Skip to main content

media_pp\elements\sink/
packet_counter.rs

1use std::sync::{
2    Arc,
3    atomic::{AtomicUsize, Ordering},
4};
5
6use crate::pp_log::{PpLog, pp_info};
7
8use crate::{
9    buffer::MediaBuffer,
10    control::ControlMsg,
11    element::{Element, ElementType, Sink, element_pp_log},
12    error::Result,
13};
14
15/// Terminal sink that just counts packets. Backed by an `Arc<AtomicUsize>`
16/// so the count can be read from outside the pipeline even when this sink
17/// ends up running on a `Queue` worker thread.
18pub struct PacketCounter {
19    pp_log: PpLog,
20    name: Arc<str>,
21    count: Arc<AtomicUsize>,
22}
23
24impl PacketCounter {
25    pub fn new(name: impl Into<String>) -> (Self, Arc<AtomicUsize>) {
26        let count = Arc::new(AtomicUsize::new(0));
27        let name: Arc<str> = name.into().into();
28        let pp_log = element_pp_log(ElementType::PacketCounter, &name, None);
29        pp_info!(pp_log: &pp_log, "created");
30        (
31            Self {
32                name,
33                pp_log,
34                count: count.clone(),
35            },
36            count,
37        )
38    }
39}
40
41impl Element for PacketCounter {
42    fn name(&self) -> Arc<str> {
43        self.name.clone()
44    }
45
46    fn element_type(&self) -> ElementType {
47        ElementType::PacketCounter
48    }
49
50    fn pp_log(&self) -> &PpLog {
51        &self.pp_log
52    }
53
54    fn pp_log_mut(&mut self) -> &mut PpLog {
55        &mut self.pp_log
56    }
57}
58
59impl Sink for PacketCounter {
60    fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
61        if let MediaBuffer::Packet(_) = buf {
62            self.count.fetch_add(1, Ordering::Relaxed);
63        }
64        Ok(())
65    }
66
67    fn control(&mut self, _msg: ControlMsg) -> Result<()> {
68        // Terminal, nothing to flush or forward.
69        Ok(())
70    }
71}